Repository navigation
AIP-104: Task Iteration - #62922
AIP-104: Task Iteration#62922dabla wants to merge 42 commits into
Conversation
d8a30b9 to
edad5de
Compare
There was a problem hiding this comment.
Thanks for working on this — excited to see DTI taking shape for Airflow 3.2. I've gone through the full diff and have feedback on the implementation, some are bugs that would crash at runtime, others are design choices worth iterating on.
A few high-level things:
-
No tests. ~700 lines of new production code with zero test coverage. We need tests for
IterableOperator,TaskExecutor,MappedTaskInstance,HybridExecutor,XComIterable,DecoratedDeferredAsyncOperator, and theiterate/iterate_kwargsmethods — covering success, failure, retry, deferral, and edge cases. -
Worker resilience. Since DTI runs N sub-tasks inside a single worker process, we need to think through what happens when that worker dies mid-execution — the scheduler has no record of which sub-tasks completed. Worth documenting the expected behavior and trade-offs here (and whether we want to add checkpointing later).
-
Thread safety. Several shared mutable structures (
contextdict,os.environ) are accessed concurrently from multiple threads without synchronization. This needs to be addressed before merge.
Inline comments below with specifics.
Thanks for pointing this out. As mentioned earlier on Slack, this PR is currently intended as an initial draft to demonstrate the concept and gather early architectural feedback. I agree that proper test coverage is essential before this can move forward. The plan is to add unit tests covering the components you mentioned (IterableOperator, TaskExecutor, MappedTaskInstance, HybridExecutor, XComIterable, DecoratedDeferredAsyncOperator, and the iterate/iterate_kwargs APIs), including scenarios for success, retries, failures, deferral, and edge cases. Once we converge on the architectural direction, I will add the corresponding test suite.
I agree this is an important architectural concern and worth discussing further. The goal of this prototype is to explore a trade-off between observability and scheduling overhead, @ashb and @potiuk mentioned the same remark before. If we try to preserve the same visibility and lifecycle guarantees as Dynamic Task Mapping, we essentially end up re-implementing DTM semantics, which brings back the same scheduler overhead that this approach is trying to avoid. This proposal intentionally explores a different point in that trade-off space: executing iterations within a single task while allowing controlled parallelism. That does mean the scheduler has indeed less visibility (but also less load) into the internal execution units.
Good point — thread safety needs to be handled carefully here. Regarding the task context, my understanding is that operators already receive a per-task context instance, but you're right that when running iterations concurrently we should avoid sharing mutable structures across threads. One possible approach would be to create a shallow or deep copy of the context for each execution unit to ensure isolation. If you have concerns about specific structures (e.g., os.environ or others), I'm happy to address them and introduce appropriate synchronization or isolation mechanisms where needed. |
960438c to
765fcfb
Compare
|
@uranusjr You should also review this PR since it touches several important modules :) |
b11f852 to
9f2c750
Compare
16ec1fc to
3242037
Compare
…ync item Every sync item with an execution_timeout went through _run_execute_callable in its worker thread, which sent SetExecutionTimeout to the supervisor again and tried TimeoutPosix, which only works on the main thread: the supervisor's hard-kill deadline drifted to "timeout after the last item started", up to a full timeout late, and the TimeoutPosix warning was logged once per item. With 4 sync items and a 30 s timeout: 4 sends, 4 warnings. _run_execute_callable takes enforce_timeout, False for an IndexedTaskInstance: the items run in worker threads under the parent's own limit, which the parent enforces and reported once. The async items already handled their limit themselves. A test runs three sync items with a 30 s timeout through the parent's _run_execute_callable and checks a single SetExecutionTimeout; it sees four before the change.
The callbacks section of the docs page and the class docstring describe which callbacks run per item; neither said what happens to listeners. They fire once, for the task instance, when the runner reports its state, as for any task: an item is not a task instance and fires none, where .expand() fires them once per mapped task instance.
IndexedTaskInstance.xcom_push writes under <key>_<index>, so iterations never
overwrite each other, but a pull had no counterpart since the override that
suffixed every pull was removed: ti.xcom_pull(key="progress") inside an item
read the parent's unsuffixed key and got None for the value the item had just
pushed, and reading it back needed key=f"progress_{ti.index}". The task state
store already follows the symmetric rule: its accessor suffixes get, set and
delete alike.
xcom_pull and axcom_pull now add the index for a pull of the iteration's own
XComs, when no task is named or its own task in its own DAG is, and leave a
pull from another task, from several, or from another DAG untouched, which is
what the review of 2026-09-10 asked for. The docs page and the class docstring
say so, and how to read another iteration's value.
Tests fail before: which key reaches the base method for six kinds of pull
and for the async pull, and an item that pushes a key and reads it back
through both forms of its own pull.
_run_tasks told apart, inline, the item outcomes the iteration cannot carry (a deferral, a reschedule, a DAG run trigger, a downstream skip, a BaseException that is no Exception) from the failures it collects, in a chain of isinstance checks that grew with every round. They move into IterableOperator._fail_fast_for, which returns the AirflowFailException that ends the task at once for such an outcome, or None; the loop raises what it returns and collects the rest. Same messages, same behaviour; the operator's own policy stays on the operator rather than on the expand input, since the exceptions come from running the items, not from resolving the input.
The exception handling of the iteration was spread over four protected methods of IterableOperator and a 95-line _run_tasks, with the failed runners and the kept retry policy decision living on IterationState. IndexedTaskOutcomes now holds all of it: a context manager entered next to Checkpoints around the loop, whose record() takes each indexed task's outcome (count, skip, collect, or raise at once for an outcome the iteration cannot carry) and whose conclude() raises the task's own outcome from what was recorded and from the run's IterationState (killed, failed on the exception the runner judges, empty input, every indexed task skipped). Whatever exception leaves the block, the failed indexed tasks' callbacks are reported on exit with the task's fate, so a fail-fast raised from record() and the outcomes raised by conclude() take the same path. IterationState keeps only what the kill needs. The prose of both modules says "indexed task" where it said "item", and _run_task takes the run's collector as a keyword argument.
The module-level functions of iterableoperator.py each served one class. The constructor checks and the input fingerprint pair are now methods of IterableOperator; the outlet-event snapshot and its replay belong to the checkpoint, IndexedTaskState, which owns that field; the live merge of one indexed task's events into the parent's belongs to IndexedTaskRunner, which owns the accessors. The replay writes into the parent's accessors directly instead of building a throwaway set and merging it. The retry decision weights are an attribute of IndexedTaskOutcomes, their only user.
…ill that arrives before the run IndexedTaskStateStoreAccessor wraps the parent's accessor without calling its constructor, so the inherited _clear_backend_only read a _scope the view never had; it now refuses as clear() does. The run's IterationState is no longer replaced when _run_tasks starts, which discarded an on_kill() that landed before the run: execute() renews it once the run ended, so a kill before the run stops it and a rerun in the same process starts clean. XComIterable.deserialize returns cls(**data); a conftest comment reads again.
The success and skip callbacks fired from IndexedTaskRunner.__exit__, as soon as the operator returned and before the indexed task's checkpoint and result were written. A checkpoint write that failed had already announced a success, and the retry ran the indexed task again and announced it a second time; the callback also found no end_date on the task instance. A plain task fires its callback after the push and the state report, with end_date set. The exit now only notes a failure. IterableOperator._run_task calls the runner's report_success() or report_skip() once the matching checkpoint is written, where the indexed task's code ran (its worker thread for a sync operator, the loop for an async one); both set end_date and the state before the callback. A result push that fails after the checkpoint is replayed from it, so the callback fires once. The docs page and the class docstring say so. test_the_success_callback_waits_for_the_checkpoint and the two runner report tests fail before.
…valuating the policy again With several failed indexed tasks whose retry policy decisions were all the default, or whose evaluation raised, no exception outweighed the others, the group was handed to the runner and no decision was kept for it. The callbacks then evaluated the policy once more, on the BaseExceptionGroup rather than on any indexed task's exception: K+1 evaluations instead of K, and for a policy that answers differently per call a retry announced here that the runner's own evaluation of the group could refuse. The once-per-item test did not reach this path because its policy answers retry(), so an exception was always chosen. _failure_for_the_runner now keeps the default decision for the group when no decision outweighs it, so _task_will_retry follows it and falls back to retry eligibility without evaluating the policy again. test_an_undecided_group_is_not_evaluated_again_for_the_callbacks fails before with four evaluations.
…the iteration on_kill() reaches the sub-operators that have started, and the stop flag keeps the executor from pulling more. An indexed task instance is neither until its runner starts: on a retry attempt it first reads its checkpoint, and a kill that lands during that read found it unregistered, so it started afterwards, ran to completion unkilled, its remote work included, and the drain waited for it before the task could conclude as terminated. IterableOperator._run_task now asks the iteration state whether a stop was requested once the checkpoint replay is done and before the runner is built, and returns IndexedTaskInstanceNotStarted as the outcome: no code run, no checkpoint, no callback, and the next attempt runs it. IndexedTaskOutcomes counts these apart from the indexed tasks that ran and the kill message names them. A sync indexed task is still handed to the pool afterwards, whose threads are as many as the calls in flight, so only that pickup remains. The docs page and the class docstring list it with the other outcomes that fire no callback. test_an_item_pulled_before_the_kill_but_not_started_does_not_run fails before with two of three indexed tasks run; the outcomes unit test covers the count.
The static checks fail on D401 for IndexedTaskInstance.context_for: its summary line named what the method returns instead of saying what it does.
The message claimed "the rest never pulled" also when every item had been pulled, and said the input was being resolved for a kill that landed before the run started, where nothing was. IndexedTaskOutcomes._killed_message now names the items that ran, those pulled before the kill and never started and those never pulled, each only when the count is not zero, and a kill before the input was resolved says so. test_a_kill_after_every_item_was_pulled_claims_no_remainder fails before.
…puts compare by it
…nt's execution_timeout
… would leave taken
… task instance's identity
Was generative AI tooling used to co-author this PR?
Claude Code (Fable 5.1).
Generated-by: Claude Code (Fable 5.1) following the guidelines
Description
This PR is the initial implementation of Iterable Tasks (IT), as discussed in the devlist and building upon the foundations of AIP-104. (Originally prototyped as "Dynamic Task Iteration"; renamed to Iterable Tasks following review feedback to avoid confusion with Dynamic Task Mapping.)
For further context on the use cases and performance benefits of IT, see this Medium Article and the new
dynamic-task-mapping-vs-iteration.rstdoc added in this PR, which compares IT with Dynamic Task Mapping (DTM) and Dynamic Task Batching in depth.The XCom Database Constraint Challenge
While porting our internal "monkey-patched" version of IT (used since Airflow 2.x) to the core, I've identified a significant technical hurdle regarding XCom handling.
Around Airflow 2.10/2.11, a change was introduced to the database constraints for the XCom table. Specifically:
map_index >= 0) unless a corresponding mapped TaskInstance exists in thetask_instancetable.The drawback is that XComs wouldn't automatically be removed from the database when a TaskInstance is deleted, which is the purpose of that constraint. So appending the index to the XCom key would be a good enough solution for IT, but not for DTM.
Current implementation in this PR
XComIterable(airflow.sdk.bases.xcom) appends the sub-task index directly to the XCom key (return_value_<index>) to bypass the constraint, and exposes the results as a lazySequence(__len__/__getitem__/__iter__) so a downstream task can consume them the same way it would consume an.expand()result.IterableOperatortracks per-sub-task progress in the AIP-103 Task State Store rather than XCom, and participates in Airflow's standard retry mechanism: if the IterableOperator's task instance is retried (or manually cleared before it finishes), already-succeeded sub-tasks are skipped and only pending/failed ones re-run. XCom is only (re)written per index once a sub-task succeeds; once every index has succeeded, a completion marker is written so that a subsequent manual clear reruns all indices from scratch rather than replaying stale results. The checkpoints themselves are kept for the store's retention ([state_store] default_retention_days), so a manual clear of a task that exhausted its retries resumes from them.on_killpropagation are handled per sub-task: a checkpointed sub-task replays its recorded outlet events on the following attempt instead of losing them, and killing the IterableOperator's task instance propagates to any sub-tasks still in flight.I believe the cleanest long-term path is still to add a dedicated route in the Execution API that retrieves multiple XComs for a single TaskInstance by a list of keys in one round trip, so
XComIterable.__getitem__/slicing don't need one request per element. I have a PR open to address this, intentionally split out of this PR.This was also discussed in the devcall, see 2026-06-04 Dev Call Minutes.
AIP-104 itself was discussed again in the latest devcall, where the concerns raised there have also been addressed: 2026-09-10 Dev call Minutes.
Examples
The examples below assume an HTTP connection named
pokeapipointing tohttps://pokeapi.co.Task Iteration
This example fetches a list of Pokémon from the PokéAPI and then uses Iterable Tasks (IT) to retrieve the details of each Pokémon. A single task instance processes all Pokémon URLs.
Comparison
get_pokemon.expand(url=urls)get_pokemon.iterate(url=urls)This demonstrates how Task Iteration can significantly reduce TaskInstance creation overhead. Task Spreading (running one iteration over exactly
NTaskInstances with.spread(across=N).iterate()) is split out into #73688.Notable design points addressed since the initial draft
.iterate()'s dict-argument semantics now match.expand(): passing adictvalue forwards(key, value)pairs to each sub-task instead of bare keys.IterableOperator.operator_nameforwards to the wrapped operator (including@task-decorated callables with acustom_operator_name), so sub-tasks report the correct type in the UI/API instead of always showingMappedOperator/IterableOperator.XComIterable.flatten()moved out to Add XComIterable.flatten() to read an iterated task's pages as one sequence #73807, stacked on this PR, so this PR stays about running a task over its input..iterate()refuses a literal that is no collection (a scalar, a string,None) at parse time on both the classic and the@taskpath, as.expand()does, and the same value from an upstream at run time.execution_timeoutcaps the whole iteration; no limit applies per indexed task, sync or async (one of the same value, started later, could never fire first). AnAirflowTaskTimeouta sync indexed task raises itself (a hook that gave up waiting) is that task's own failure, as it is the mapped task instance's under.expand().on_kill()takes no lock, so the runner's SIGTERM handler cannot hang on the loop thread's own register/unregister; the copyprepare_for_executionmakes gets an iteration state of its own, so a kill in one attempt does not stop the next underdag.test(); anasendcancelled while waiting for the comms thread lock no longer leaves it taken (CommsDecoder._acquire_thread_lock).KubernetesPodOperatorreattaching,EcsRunTaskOperator(reattach=True)) are listed under "When not to use IT", with{{ ti.index }}as the label that tells items apart.on_kill()propagates to in-flight sub-tasks, stops the iteration from starting any other (one pulled before the kill but not yet started is counted apart and left for the next attempt), and outlet/asset events recorded by a sub-task that already succeeded are replayed from its checkpoint on a later retry instead of being lost.TriggerDagRunOperator, andShortCircuitOperator-style downstream skipping are explicitly rejected insideIterableOperatorwith actionable errors rather than being silently mishandled — see the class docstring for the full list of current limitations.multiple_outputsis ignored for iterated tasks, explicitly at theIterableOperatorlevel. A@taskwith aMappingreturn annotation infersmultiple_outputs=True, but the value the runner pushes for an iterated task is theXComIterableaggregate rather than a dict, so honouring the flag made the runner reject the result after every sub-task had already succeeded. Each sub-task's return value is pushed whole asreturn_value_<index>; keys are not fanned out into separate XComs the way.expand()does. Documented on the class and in the Task SDK docs, and pinned by a runner-level regression test for a dict-returning task under.iterate().Per-iteration keys: XComs and task state
Every iteration of an iterated task runs under the same task instance (same dag id, task id, run id and map index). Anything an iteration writes into a per-task-instance store therefore competes with its siblings for the same key, and with the async executor the winner is whichever iteration finishes last. Two stores are affected, and both now apply the same rule: a key written from inside an iteration carries that iteration's index.
IndexedTaskInstance.xcom_push/axcom_pushsuffix the key with_<index>. That is what makesreturn_value_<index>andXComIterablework, and it applies to any key an operator pushes fromexecute, including the keys of amultiple_outputsdict. Pulls are not suffixed:ti.xcom_pull(task_ids="upstream")reaches the upstream's XCom untouched.context["task_state_store"]). This was a gap: an iteration that stored a watermark or a cursor withtask_state_store.set("last_offset", ...)shared that key with every sibling.IndexedTaskStateStoreAccessorcloses it:IndexedTaskInstance.task_state_storeis the parent's accessor seen through the index, suffixing keys onget/set/deleteand their async twins, and the sub-task's context carries the same object, so an operator does not need to know it is being iterated.clear()is refused inside an iteration, since it would wipe the siblings' state and the operator's own checkpoints; an iteration deletes its own keys instead.IterableOperatorrecords per-index progress in the parent's store under_iterable_<index>and_iterable_completed, written through the parent's accessor, so they are never double-suffixed and never collide with user keys.IndexedTaskRunner(formerlyTaskExecutor, renamed because it read like one of Airflow's executors) builds the context an iteration runs against: a copy of the parent's context with the iteration's own task instance, its indexed state store view and its own outlet events. The operator binds it from thewithstatement, runs the operator inside that block, and records the outcome (checkpoint, XCom push, outlet-event merge) after it, soon_killand the failure callbacks apply to the operator's execution only.{pr_number}.significant.rstor{issue_number}.significant.rst, in airflow-core/newsfragments.